相信很多人在了解 Comet 後,可能會思考一個問題,Comet 夾在 Spark 與 datafusion 中間,如果遇到datafusion不支援的運算怎麼辦?會不會直接失敗呢?身為某大廠撐腰的專案,當然不可能讓這件事發生,因此今天要來介紹 Comet fallback 機制
一般以為只有兩條:
實際上還有第三條,Spark 自己 doGenCode 產出的 Java 碼被包成 Arrow batch kernel,直接塞進 Comet pipeline 跑。這條路不會換掉 Comet operator、資料不從 Arrow 轉回 UnsafeRow,依然還在 native pipeline裡。
這個是 Comet 1.0 加入的表達式級 dispatcher,核心要解決一個叫白白 fallback 的痛
先來看一句 SQL:
SELECT upper(s), next_day(d, day COLLATE utf8_lcase) FROM t
Project 中有兩個表達式,如果 Comet 想把整個 Project 換成 CometProject,條件是每一個表達式都必續轉成 protobuf,一但轉失敗就整段 Project 退回 Spark。
upper(s) Comet 有原生 Rust 實作,本來來是可以走 native,但是 next_day(d, day COLLATE utf8_lcase) 用到了 UTF8_LCASE collation,而 Rust 那一側目前沒對齊這個 collation 的行為,回報 Incompatible(collationReason)。
而這個整段就會變成 Spark Project,但 upper(s) 原本是可以走 native,也需要跟退回 Saprk,且資料還要在中間插一個 ColumnarToRow 從 Arrow 轉回 UnsafeRow。
這就是爲什麼叫做白白 fallback,因為一個表達式不能 native,全部連帶賠掉。
Comet 1.0 之後,同一個表達式的路徑選擇長這樣:
| 路徑 | 誰在算 | 資料還在 Comet 管線裡嗎 | Operator 前綴 |
|---|---|---|---|
| Native | Rust kernel(DataFusion) | 在 | CometProject |
| JVM codegen dispatcher | Spark 產生的 Java 碼,包成 batch kernel | 在(仍是 CometProject) |
CometProject |
| Fallback 回 Spark | Spark 的 Project |
不在(插入 ColumnarToRow) |
Project |
中間那條就是今天的主角,EXPLAIN 印出來會看到:
CometProject [COMET-INFO: JVM codegen dispatcher: next_day]
Operator 前綴還是 Comet*,只是特定表達式走 dispatcher,不是 Rust。
運算沒退回 Spark,只是那顆表達式的計算換人做。
每個 Spark 表達式都有 doGenCode,這個方法不是直接算值,而是吐出一段 Java 原始碼,描述「這一列該怎麼算」。這條路是 Spark 2.0 的 whole-stage codegen(WSCG)的基石,運算子把上下游的 doGenCode 拼成一段大函式,交給 Janino 編譯,跑起來就跟手寫的 Java 差不多。
Comet 這邊的做法是單獨呼叫這一個表達式的 doGenCode,把 Spark 生出來的那段 Java 包成一顆吃 Arrow batch 的 kernel。這不是「在 JVM 上跑一整棵 Spark plan」,而是「Comet 管線裡嵌一顆用 Spark 邏輯寫成的 Arrow kernel」。
好處三個:
代價是每次呼叫都要跨 JNI 進 JVM,dispatcher 的成本比純 Rust 高。所以它是選擇性的補位,不是所有表達式都該走。
決定哪些表達式該走 dispatcher,這個開關長成一個 Scala trait
定義(CometExpressionSerde.scala):
/**
* Mixin for serdes whose native implementation cannot match Spark for some inputs. When
* getSupportLevel returns Incompatible and the user has not enabled allowIncompatible, the
* expression routes through the JVM codegen dispatcher ... instead of falling the projection
* back to Spark. When getSupportLevel returns Unsupported, the expression always routes
* through the dispatcher ...
*
* Enrollment is opt-in: only serdes that explicitly mix this in are routed through the
* dispatcher. Every other Incompatible expression falls back to Spark, and every other
* Unsupported expression falls back to Spark.
*/
trait CodegenDispatchFallback extends NativeOptInAvailable { self: CometExpressionSerde[_] => }
trait 本身沒有方法要實作
契約寫在 scaladoc:混入這個 trait 等於報名「我這條 native 路走不了時,請改問 dispatcher,不要直接放棄整個投影」。
沒混入就代表「Incompatible / Unsupported 就是整段退回 Spark」
這個設計是顯性 opt-in
每一個 serde 要不要 dispatcher 補位,作者自己決定,不是 default-on。原因是 dispatcher 有 JNI 成本,如果 default-on 所有 serde 一律走它,可能出現「明明可以完全 native 卻走了 dispatcher」的意外性能坑。
實際 route 的邏輯在 QueryPlanSerde.exprToProtoInternal:
handler.getSupportLevel(expr) match {
case Unsupported(notes) =>
// `CodegenDispatchFallback` serdes have no native path for these cases either, but the
// dispatcher can still run Spark's own `doGenCode` inside the Comet pipeline. Try that
// before falling the projection back to Spark.
dispatchIfFallback(handler, expr, inputs, binding)
.orElse { /* record unsupported, fall back to Spark */ }
case Incompatible(reason) if !allowIncompatible =>
dispatchIfFallback(handler, expr, inputs, binding)
.orElse { /* record incompatible, fall back to Spark */ }
case _ => /* native path */
}
dispatchIfFallback 只看一件事:這個 serde 有沒有混入 CodegenDispatchFallback。有,就呼叫 emitJvmCodegenDispatch 生 dispatcher plan;沒有,就回 None,整個 Project 掉回 Spark。
分界很清楚:Unsupported 或 Incompatible 本身不是路徑決定,只是 handler 告訴 planner「native 這條走不了」。真正決定「還能不能救」的是有沒有那個 mixin。
回到開頭的例子:
CometNextDay 回 Incompatible(collationReason)
CometLevenshtein 回 Unsupported(...)
沒有 mixin 時,兩者都等於「轉換失敗」,整段 Project 退回 Spark。
加上 mixin 之後:
next_day(...) 走 dispatcher(Spark 的 doGenCode,連 Collate 子樹一起融進同一顆 kernel)upper(s) 繼續 nativeCometProject
「白白」指的就是這件事:dispatcher 這條路早已存在,只要混入就能用,答案跟 Spark 逐位元組相同,卻因為沒混入而整段 Project 被丟掉。
Scala 的 mixin,就是用 with 把一個 trait 混進既有的 class / object,讓它多出一份行為或標記,而不改繼承樹。以 CometNextDay 為例:
object CometNextDay extends CometExpressionSerde[NextDay] with CodegenDispatchFallback
讀法:CometNextDay 本體是 CometExpressionSerde[NextDay],同時混入 CodegenDispatchFallback。因為 CodegenDispatchFallback 沒有抽象方法,混入不需要實作任何東西,就是貼一張報名貼紙。
Java interface 通常是「你必須實作這些方法」。這裡的 trait 更像貼紙:
convert
getSupportLevel 回 Unsupported / Incompatible 時,讓 dispatchIfFallback 多試一條路執行期靠 handler match { case h: CodegenDispatchFallback => ... } 認出誰報名了。self: CometExpressionSerde[_] => 這個 self-type 限制只有 serde 才能混這個 trait,避免亂貼到別的類別上。
Java 21 之後也有類似的空 interface 加 sealed 加 pattern matching 的做法,但 Scala 的 mixin 更輕,多個 trait 可以疊。例如某個 serde 同時 with CodegenDispatchFallback 又繼承 NativeOptInAvailable,文件會標為 Hybrid serde。這是 mixin 組合,不是單一路徑繼承。
Comet 表達式的三條路:native、JVM codegen dispatcher、Spark fallback
第三條 dispatcher 把 Spark 自己 doGenCode 產出的 Java 包成 Arrow batch kernel,塞進 Comet 管線裡跑,答案跟 Spark 逐位元組相同,資料不用轉回 UnsafeRow。這條路靠一個空 trait CodegenDispatchFallback 當開關,每個 serde 顯性 opt-in,避免默默走 dispatcher 掩蓋語義漏洞。這種「空 trait 當貼紙、執行期 pattern match」是 Scala mixin 跟 Java interface 最大的差別。
明天 D22 講plan 攔截,Spark physical plan 在什麼時機被 Comet 換掉?AQE 執行期重規劃時,攔截規則怎麼保持一致而不衝突?這是 Comet 之所以能塞進 Spark 而不打破 Catalyst 的關鍵。
那就明天見~
QueryPlanSerde.scala(exprToProtoInternal route 邏輯): https://github.com/apache/datafusion-comet/blob/main/spark/src/main/scala/org/apache/comet/serde/QueryPlanSerde.scala
CometExpressionSerde.scala(CodegenDispatchFallback trait 契約): https://github.com/apache/datafusion-comet/blob/main/spark/src/main/scala/org/apache/comet/serde/CometExpressionSerde.scala
Expression.scala(doGenCode 定義): https://github.com/apache/spark/blob/master/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/Expression.scala